02 - LangGraph Checkpointer
前置:01 篇的检查点概念。
本篇回答:LangGraph 存了什么、恢复点怎么算出来、崩溃窗口怎么处理、以及它做不到什么。
术语
- 通道(channel):节点之间传数据的位置,可先理解为状态里的一个字段
- 超级步(superstep):一轮并行执行。执行模型借自 Google Pregel —— 一轮内所有待执行节点并行运行,统一提交后进入下一轮
- thread:一次会话或一个任务的标识,是检查点的主键
一、检查点的数据结构
libs/checkpoint/langgraph/checkpoint/base/__init__.py 中的定义是理解整个机制的入口:
class Checkpoint(TypedDict):
"""某一时刻的状态快照。"""
v: int
# 检查点格式版本号,当前为 1。
# 存在这个字段意味着格式会演进 —— 旧检查点必须能被新代码读懂。
id: str
# 检查点 ID。同时满足唯一且单调递增,因此可直接用于排序。
# 这样列历史时不必依赖时间戳 —— 分布式环境下时钟不可 靠
# (节点间漂移、同一毫秒内产生多条)。
ts: str
# ISO 8601 时间戳。仅供展示,不参与排序。
channel_values: dict[str, Any]
# 各通道在该时刻的实际值,即"状态本身"。
channel_versions: ChannelVersions
# 各通道「当前」的版本号。每次写入该通道,版本号递增。
versions_seen: dict[str, ChannelVersions]
# 每个节点"看过"每个通道的哪个版本。
# 与 channel_versions 比对即可算出哪些节点需要执行 —— 见 1.2 节。
updated_channels: list[str] | None
# 本次检查点更新了哪些通道,用于减少下游的比对范围。
1.1 ID 单调递增带来的简化
id 同时承担主键与排序键两个职责。这个决定省掉了一整类问题:列历史记录、找父检查点、判断先后顺序全部退化成整数比较,不需要额外的序号字段,也不受时钟影响。
1.2 恢复点是算出来的,不是存下来的
这是整个机制里最值得理解的一点:LangGraph 没有"执行到第几步"这样的字段。
它存的是版本对照关系:
调度逻辑就是一句话:某节点看过的版本落后于通道当前版本,该节点就进入待执行集合。
这个设计换来三件用"程序计数器"做不到的事:
| 能力 | 为什么成立 |
|---|---|
| 并行分支 | 多个节点各自落后于不同通道,一轮内可全部调度 |
| 人工改状态后自动续跑 | 手工写入某通道 → 版本号递增 → 依赖它的节点自动变为待执行 |
| 图结构变更不影响恢复 | 依赖关系按通道计算,与节点顺序无关 |
可迁移的结论:为 DAG 或状态机做恢复时,记录"每个节点消费到哪个版本"比记录"全局执行到第几步"更健壮 —— 前者是局部信息、天然 可并行、对拓扑变化不敏感。
1.3 全量值与增量存储
channel_values 存全量。当某个通道很大(例如一段长对话历史)时,每步存一次会造成显著的写放大 —— 即 01 篇 2.4 节描述的问题。
LangGraph 的解法是 DeltaChannel(当前标注 beta):平时只存增量写入,周期性存一次完整快照作为种子,恢复时从最近快照沿增量链向前叠加。
元数据里记录了触发快照的两个计数:
(updates, supersteps)
updates —— 有多少个超级步写过这个通道
supersteps —— 距离该通道上次快照,总共过了多少个超级步
官方注释说明了第二个计数的作用:
The supersteps bound prevents unbounded ancestor walks on threads where [the channel is rarely written]
很少被写的通道,其"最近快照"可能在几百步之前,恢复时要回溯极长的祖先链。所以除按写入次数触发快照外,还需要一个按超级步数量的强制上限兜底。
这与数据库的 WAL + checkpoint、Redis 的 AOF + RDB 是同一模式:增量存储必须配强制全量点。
二、中间写入:崩溃窗口的处理
01 篇 2.2 节指出,"步骤完成"与"检查点落盘"之间必然存在窗口。LangGraph 用 pending_writes 处理框架内部的这段窗口。
2.1 机制
class CheckpointTuple(NamedTuple):
config: RunnableConfig
checkpoint: Checkpoint
metadata: CheckpointMetadata
parent_config: RunnableConfig | None = None
pending_writes: list[PendingWrite] | None = None
# ↑ 已完成但尚未并入下一个检查点的节点输出
def put_writes(
self,
config: RunnableConfig, # 关联到哪个检查点
writes: Sequence[tuple[str, Any]], # (通道名, 值) 列表
task_id: str, # 哪个节点产生的 —— 恢复时靠它判断该节点是否已完成
task_path: str = "",
) -> None:
"""存储关联到某个检查点的中间写入。"""
方法文档里的关键词是 intermediate:一个超级步内多个节点并行执行,每个节点完成后立即单独落一次写入并带上 task_id;整轮结束后才生成新的 Checkpoint。
没有这个机制,整个超级步的三个节点都要重来。在"一个节点 = 一次模型调用"的场景下,这是三倍的直接成本。
2.2 它的作用范围
pending_writes 保证的是框架不会重复执行已完成的节点。它不能保证节点内部调用外部服务的幂等 —— 如果节点 A 已经发出飞书消息但 put_writes 写入失败,恢复后消息仍会重发。
外部副作用的幂等仍然由 01 篇第三节的幂等键负责,框架层解决不了。
三、接口全貌
BaseCheckpointSaver 的同步方法(每个都有 a 前缀的异步版本):
| 方法 | 作用 | 类别 |
|---|---|---|
get / get_tuple | 读取某个检查点 | 基础 |
list | 列出历史 | 基础 |
put | 写入检查点 | 基础 |
put_writes | 写入中间结果 | 基础 |
delete_thread | 删除某 thread 的全部检查点 | 运维 |
delete_for_runs | 按 run 删除 | 运维 |
copy_thread |